Skip to content

test: fix flaky EmbeddedKafkaSupervisorTest header filtering test - #20519

Merged
kfaraz merged 1 commit into
apache:masterfrom
confluentinc:sugoel/fix-flaky-kafka-header-filtering-test
Oct 8, 2026
Merged

kfaraz merged 1 commit into
apache:masterfrom
confluentinc:sugoel/fix-flaky-kafka-header-filtering-test

Conversation

@suraj-goel

@suraj-goel suraj-goel commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #20518.

Description

EmbeddedKafkaSupervisorTest.test_runKafkaSupervisorWithHeaderFiltering (added in #18525) failed its first attempt in 6 of 7 master runs. Each time, it timed out waiting for ingest/handoff/count to reach 4. It fails only when it runs after the other tests in the class. The surefire retry runs it alone, so the retry passed.

Root cause

The task does read all 10 records. The race is between the test suspending the supervisor and the task's own checkpoint:

  1. With maxRowsPerSegment = 1, the task requests a checkpoint right after it reads its batch.
  2. The test suspended the supervisor as soon as the broker emitted any metric for the datasource. In a warm cluster, that happens within milliseconds of the first segment allocation. The suspend therefore replaces the supervisor while the checkpoint request is in flight.
  3. The overlord rejects the checkpoint (supervisor could not be found). The task fails with Checkpoint request ... failed, dying, and nothing is published.

The spec also used the default startDelay (PT5S) and period (PT30S). The supervisor adds a newly created task to its task group only on its next run. Until then, a checkpoint request does nothing (New checkpoint is null) and leaves the task paused. The test depended on the suspend for publishing: the suspended supervisor's first run, 5 seconds later, set the end offsets.

The fix suggested in the issue (wait for ingest/events/processed >= 4 before suspending) makes the race less likely but does not remove it. processed is incremented during parsing, a few milliseconds before the checkpoint request is sent, and the metric is emitted every 100 ms.

Fix

  • The header-filter supervisor spec is now built from newKafkaSupervisor(), like the other tests in the class, with only the header filter added. The removed builder was identical apart from the timing settings: 500 ms task duration, 500 ms run period, and 10 ms start delay.
  • The test now waits for handoff without suspending the supervisor. The short task duration publishes the segments. The existing @AfterEach teardown then suspends the supervisor and cancels the remaining tasks.

Suspending a supervisor while a task's checkpoint is in flight fails that task. The failure is recoverable: the task's offsets were not committed, so resumed tasks re-read the data. This PR does not change that behavior.

Testing

Ran locally on JDK 25 with -Dsurefire.rerunFailingTestsCount=0:

Version Full class Test alone
Before 0/4 passed (same timeout at line 220 as CI) 5/5 passed
Issue's suggested fix 6/6 passed (race still possible) not run
This PR 8/8 passed 3/3 passed

With this change, handoff completes about 1.6 seconds after the supervisor starts, compared with about 5.5 seconds after the suspend before.


Key changed/added classes in this PR
  • EmbeddedKafkaSupervisorTest

This PR has:

  • been self-reviewed.
  • added comments explaining the "why" and the intent of the code wherever would not be obvious for an unfamiliar reader.
  • added unit tests or modified existing tests to cover new code paths, ensuring the threshold for code coverage is met.

@kfaraz kfaraz left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, @suraj-goel !
Before merging off this PR, could you also confirm if this change has affected the runtime of this test in anyway? Just curious.

Edit: never mind, read the PR description. It seems like the test will now finish slightly faster?

@suraj-goel

Copy link
Copy Markdown
Contributor Author

Thanks, @suraj-goel ! Before merging off this PR, could you also confirm if this change has affected the runtime of this test in anyway? Just curious.

Thanks @kfaraz! Yes, it's faster. Measured locally:

  • Running only this test (including embedded cluster startup): ~23.4s → ~14.4s
  • Supervisor submit → last segment handoff: ~10.7s → ~1.75s

@suraj-goel

Copy link
Copy Markdown
Contributor Author

@kfaraz , I don't have permissions to merge. Feel free to merge this PR.

@kfaraz
kfaraz merged commit cde7404 into apache:master Oct 8, 2026
28 checks passed
@github-actions github-actions Bot added this to the 39.0.0 milestone Oct 8, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Flaky test: EmbeddedKafkaSupervisorTest.test_runKafkaSupervisorWithHeaderFiltering

2 participants